Skip to content

Distinguish OpenAI batch timeout and cancellation in deferrable tasks - #72149

Merged
Lee-W merged 8 commits into
apache:mainfrom
astronomer:openai-batch-termination-reason
Sep 29, 2026
Merged

Lee-W merged 8 commits into
apache:mainfrom
astronomer:openai-batch-termination-reason

Conversation

@Lee-W

@Lee-W Lee-W commented Aug 27, 2026 •

Copy link
Copy Markdown
Member

A deferred OpenAITriggerBatchOperator reported every non-success outcome as the same OpenAIBatchJobException, so a task that timed out looked identical to a batch that failed, and neither could be handled separately. The synchronous path has always raised OpenAIBatchTimeout for the same condition.

The trigger now reports why it stopped, and the resuming task picks the matching exception from that field rather than from the message text. When a deferred task times out, the operator also requests cancellation of the batch, as the synchronous path already did, instead of leaving it running on OpenAI's side.

Cancellation is best-effort everywhere it is requested (deferred timeout, synchronous timeout, on_kill): a failed cancel request is logged as a warning and no longer replaces the exception that explains why the task stopped.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: [Claude] following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch 3 times, most recently from 57ed517 to fc9b2fc Compare September 17, 2026 06:21
@Lee-W
Lee-W marked this pull request as ready for review September 17, 2026 06:39
@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch from 8b92065 to 33179d9 Compare September 18, 2026 09:27
Comment thread providers/openai/src/airflow/providers/openai/operators/openai.py Outdated
Comment thread providers/openai/src/airflow/providers/openai/operators/openai.py
Comment thread providers/openai/src/airflow/providers/openai/hooks/openai.py Outdated
@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch from 33179d9 to 934d00c Compare September 19, 2026 01:27
@Lee-W
Lee-W marked this pull request as draft September 19, 2026 07:54
@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch 2 times, most recently from 1763aca to 44bcb7e Compare September 24, 2026 03:00
@Lee-W
Lee-W marked this pull request as ready for review September 24, 2026 07:26
@Lee-W
Lee-W requested a review from kaxil September 24, 2026 07:50
@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch from 44bcb7e to 4aef57e Compare September 28, 2026 02:55
Comment thread providers/openai/docs/changelog.rst Outdated
Comment thread providers/openai/src/airflow/providers/openai/operators/openai.py Outdated
Comment thread providers/openai/tests/unit/openai/operators/test_openai.py Outdated
Comment thread providers/openai/docs/changelog.rst Outdated
Comment thread providers/openai/src/airflow/providers/openai/operators/openai.py Outdated
Comment thread providers/openai/tests/unit/openai/operators/test_openai.py Outdated
A deferred OpenAITriggerBatchOperator reported every non-success outcome as the
same OpenAIBatchJobException, so a task that merely timed out looked
identical to a batch that failed, and neither could be handled separately. The
synchronous path has always raised OpenAIBatchTimeout for the same condition.
The trigger now reports why it stopped, and the resuming task picks the matching
exception from that field rather than from the message text.
The synchronous path cancels the batch before raising a timeout, but the
deferrable path only failed the task and left the batch running and billing on
OpenAI's side. A deferred timeout now requests cancellation with the batch id
carried by the trigger event. Cancellation is asynchronous, so the batch
reports "cancelling" for a while before it settles; a cancel that fails is
logged and never masks the timeout.

Deliberately not cancelling on a polling_error: batches.cancel is
irreversible, and the failure mode there is unknown and often a transient,
Airflow-side error rather than a real batch problem. OpenAI's batch has a
24-hour completion window that bounds it on its own, so leaving it alone costs
a bounded amount; wrongly cancelling it is unrecoverable data loss. Anthropic
cancels on its equivalent error branch, but its session has no such natural
end point, so that precedent does not apply here.
Every TriggerEvent OpenAIBatchTrigger yields sets batch_id from a required
constructor argument, and OpenAITriggerBatchOperator only ever defers to that
trigger, so the branch that logged "carried no batch_id" could not be reached
and had no test covering it. The batch id is now read the way the success path
in the same method already reads it, with event["batch_id"], so the two paths
are consistently defensive.
OpenAITriggerBatchOperator.on_kill and OpenAIHook.wait_for_batch's timeout
branch both called cancel_batch bare, so an error from the cancellation request
propagated in place of the reason the task was actually stopping. on_kill now
goes through _cancel_batch_quietly, and the synchronous timeout path logs a
warning and still raises OpenAIBatchTimeout. Both paths now behave the way the
deferred timeout path already did.
The trigger wrote termination_reason as a bare string and the hook keyed its
exception map on bare strings, so a typo on either side type-checked fine and
only a test could catch it. Both sides, and the operator's timeout check, now
go through TerminationReason, making a misspelled reason an attribute error.
The enum subclasses str, so the value crossing the triggerer boundary is
unchanged and the existing tests still pin it.
The note that a deferred timeout now raises OpenAIBatchTimeout was lost in a
rebase. It is back as a warning that spells out OpenAIBatchTimeout is not an
OpenAIBatchJobException, so handlers keyed on the latter stop matching. Both
notes from this change now sit above the 2.0.0 heading, since 2.0.0rc1 was
cut without it.
The timeout case of the reason-to-exception test used a real OpenAIHook, so
it passed only because the failed connection lookup was swallowed. It is now
one table of reason, expected exception and whether cancellation is
requested, with a spec'd hook mock for every row.
…cstrings

The timeout param now keeps to what it bounds plus the execution_timeout
caveat, scoped to deferrable mode. The termination-reason-to-exception
mapping lives only in execute_complete, and the asynchronous cancellation
caveat only in _cancel_batch_quietly.
@Lee-W
Lee-W force-pushed the openai-batch-termination-reason branch from 4aef57e to 2f65d33 Compare September 29, 2026 05:50
@Lee-W
Lee-W merged commit 4b66065 into apache:main Sep 29, 2026
79 checks passed
@Lee-W
Lee-W deleted the openai-batch-termination-reason branch September 29, 2026 08:09
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants